--- title: "01-Kafka 消费与文本内容解析" created: 2026-05-18 aliases: - Kafka 消费与文本内容解析 tags: - 项目 --- # Kafka 消费与文本内容解析 上一篇讲完了文档上传的同步链路——文件校验、MinIO 上传、数据库入库,最后把消息丢进了 Kafka。 那 Kafka 消费方拿到消息之后到底做了什么? 这篇就来拆这条异步链路的前半段:从消费消息开始,到文本提取与清洗完成为止。 先上总览流程图。 ## 异步解析链路总览 ## ![[Fhrq35ZYem8aVEVCLc9ctrbtIP-Z-b0ba0c72.png]] ## Kafka 消费者:消息入口 消费端的入口在 `DocumentKafkaConsumer`,它的职责很单一——把 JSON 消息反序列化,然后转交给异步处理服务。 ```java /** * 消费"解析路由"消息。 *
* 这一步是上传完成后的异步链入口:收到消息后,会把 documentId 和 taskId 交给异步处理服务, * 继续执行文档下载、正文解析、结构节点生成、策略推荐等步骤。 *
*/ @KafkaListener( topics = SPRING_INJECT_PREFIX_DISTINCTION_NAME + "-" + "${app.manage.kafka.parse-topic}", groupId = "${app.manage.kafka.group-id}-parse") public void consumeParseRoute(String payload) { try { // 先把 JSON 还原成强类型消息对象,避免后续处理层直接面对原始字符串。 DocumentParseRouteMessage message = objectMapper.readValue(payload, DocumentParseRouteMessage.class); // 真正的业务推进放到异步处理服务中,这里只承担"消费并转发"的职责。 asyncProcessService.handleParseRoute(message.getDocumentId(), message.getTaskId()); } catch (Exception exception) { // 消费失败只记录日志,不让异常继续向外冒泡破坏监听线程。 log.error("消费解析路由消息失败,payload={}", payload, exception); } } ``` 几个要点: - topic 名称是 `环境前缀-配置的parseTopic`,和上传端发送时用的是同一个 topic - `groupId` 带了 `-parse` 后缀,和索引构建的消费组区分开 - 整个方法用 try-catch 包住,消费失败只打日志不抛异常——这是为了防止一条坏消息把整个消费线程搞挂 - Consumer 本身不做任何业务逻辑,纯粹是"反序列化 + 转发" ## handleParseRoute:异步解析主方法 真正干活的是 `DocumentAsyncProcessServiceImpl.handleParseRoute()`。这个方法比较长,我们按阶段拆开来看。 ### 加载上下文 + 推进任务状态 ```java /** * 处理"解析路由"任务。 ** 这是文档上传成功后的第一条异步业务主链,整体顺序如下: * 1. 读取文档记录和任务记录,确认异步处理上下文存在; * 2. 把任务状态切到 RUNNING,把当前阶段推进到 CONTENT_PARSE; * 3. 从对象存储下载原始文件,并调用解析器提取纯文本和结构信息; * 4. 把解析后的纯文本重新上传为 txt,便于后续索引构建直接复用; * 5. 用结构节点服务替换文档结构节点,并同步导航索引、图投影和画像; * 6. 基于解析结果调用策略服务生成推荐切块方案; * 7. 把推荐方案和步骤写入数据库,同时更新文档的解析状态、策略状态和统计信息; * 8. 以成功或失败状态收尾任务,并记录任务日志。 *
** 这个方法本身不直接执行向量化和 chunk 落库,那是后续"索引构建任务"的职责; * 当前阶段的目标是把文档从"原始文件"推进到"解析完成并拿到推荐策略"。 *
*/ public void handleParseRoute(Long documentId, Long taskId) { SuperAgentDocument document = documentMapper.selectById(documentId); SuperAgentDocumentTask task = taskMapper.selectById(taskId); if (document == null || task == null) { log.warn("解析任务对应的文档或任务不存在,documentId={}, taskId={}", documentId, taskId); return; } Date startTime = new Date(); try { // 进入异步解析后,先把任务状态切成 RUNNING,当前阶段设为"内容解析"。 task.setTaskStatus(DocumentTaskStatusEnum.RUNNING.getCode()); task.setCurrentStage(DocumentTaskStageEnum.CONTENT_PARSE.getCode()); task.setStartTime(startTime); taskMapper.updateById(task); // 文档主表的解析状态也同步切到 PARSING,方便前端查询时看到实时状态。 document.setParseStatus(DocumentParseStatusEnum.PARSING.getCode()); documentMapper.updateById(document); ``` 第一件事就是从数据库把文档记录和任务记录捞出来。如果查不到(比如上传事务回滚了、或者数据被删了),直接 return,不做任何处理。 然后把任务状态从 `NEW` 切到 `RUNNING`,当前阶段设为 `CONTENT_PARSE`。文档主表的 `parseStatus` 也同步更新为 `PARSING`。这样前端轮询文档状态时就能看到"正在解析中"。 ### 记录开始日志 + 下载原始文件 ```java // 先记一条开始日志,把原始对象名落下来,便于后续排查下载与解析问题。 taskLogService.saveLog(taskId, documentId, DocumentTaskStageEnum.CONTENT_PARSE.getCode(), DocumentTaskEventTypeEnum.START.getCode(), DocumentLogLevelEnum.INFO.getCode(), DocumentOperatorTypeEnum.SYSTEM.getCode(), null, "开始解析文档内容。", Map.of("objectName", document.getObjectName())); // 重新从对象存储下载原始文件,确保异步链处理的是"上传后实际持久化成功"的那份文件。 byte[] fileBytes = storageService.downloadObject(document.getObjectName()); ``` > 为什么要重新下载文件? > > 上传阶段的 `byte[]` 是在 HTTP 请求线程里的内存数据,请求结束后就释放了。异步消费方运行在另一个线程(甚至另一个进程),所以必须从 MinIO 重新下载。这也顺便验证了文件确实已经成功持久化到对象存储。 ### 调用 Tika 解析文档 ```text // parserService 会负责提取纯文本、结构节点候选、字符数、token 数、质量等级等分析结果。 DocumentAnalysisResult analysisResult = parserService.parse(fileBytes, document.getOriginalFileName(), document.getMimeType(), DocumentFileTypeEnum.getRc(document.getFileType())); ``` 这一步是整条链路的核心——把原始二进制文件变成结构化的分析结果。我们深入看看 `TikaDocumentParserService.parse()` 的实现。 ## TikaDocumentParserService:文档解析器 ### parse 方法总览 ```java /** * 解析文档原始字节并生成标准化分析结果。 ** 这是上传后异步解析链里的核心一步,整体顺序如下: * 1. 先根据文件类型提取原始文本; * 2. 对文本做换行、空白字符和乱码清洗; * 3. 提取结构节点候选,并估算标题数量; * 4. 统计段落、字符数、token 数; * 5. 评估文档结构等级和内容质量等级; * 6. 把这些结果打包成 DocumentAnalysisResult 返回给上游。 *
*/ public DocumentAnalysisResult parse(byte[] bytes, String originalFileName, String mimeType, DocumentFileTypeEnum fileType) { // 第一步先把不同格式文件转换成统一的原始文本。 String rawText = extractRawText(bytes, originalFileName, mimeType, fileType); // 再做文本清洗,为后续标题识别、段落切分和结构抽取提供更稳定的输入。 String cleanedText = cleanupText(rawText); // 结构节点由专门的提取器负责抽取,这些节点后面会参与导航和切块策略判断。 List* Office/PDF/HTML 等复杂格式优先走 Tika; * 纯文本格式则直接按 UTF-8 解码。 * 如果 Tika 失败但 MIME/文件类型仍明显是文本,则退回到直接字符串解码作为兜底。 *
*/ private String extractRawText(byte[] bytes, String originalFileName, String mimeType, DocumentFileTypeEnum fileType) { try { if (fileType == DocumentFileTypeEnum.PDF || fileType == DocumentFileTypeEnum.DOC || fileType == DocumentFileTypeEnum.DOCX) { return tika.parseToString(new ByteArrayInputStream(bytes)); } if (fileType == DocumentFileTypeEnum.TXT || fileType == DocumentFileTypeEnum.MD) { return new String(bytes, StandardCharsets.UTF_8); } if (fileType == DocumentFileTypeEnum.HTML) { return tika.parseToString(new ByteArrayInputStream(bytes)); } return tika.parseToString(new ByteArrayInputStream(bytes)); } catch (Exception exception) { // 某些文本文件虽然 Tika 解析失败,但原始字节本身仍然可以直接按文本解码读取。 if (fileType == DocumentFileTypeEnum.TXT || fileType == DocumentFileTypeEnum.MD || mimeType != null && mimeType.startsWith("text/")) { return new String(bytes, StandardCharsets.UTF_8); } throw new IllegalStateException("Tika 解析失败: " + exception.getMessage(), exception); } } ``` 提取策略很清晰: - **PDF / DOC / DOCX / HTML**:走 Apache Tika,它能处理各种复杂格式 - **TXT / MD**:直接 UTF-8 解码,不需要 Tika - **兜底**:如果 Tika 解析失败,但文件类型明显是文本类的,就退化成直接字符串解码 ### cleanupText:文本清洗 ```java private String cleanupText(String rawText) { if (StrUtil.isBlank(rawText)) { return ""; } String cleaned = rawText .replace("\r\n", "\n") // Windows 换行统一成 Unix 换行 .replace('\r', '\n') // 老 Mac 换行也统一 .replace('\u0000', ' ') // 空字符替换成空格 .replaceAll("[\\t\\x0B\\f]+", " ") // 制表符等控制字符替换成空格 .replaceAll("\\n{3,}", "\n\n") // 连续 3 个以上换行压缩成 2 个 .replaceAll("[ ]{2,}", " ") // 连续多个空格压缩成 1 个 .trim(); return cleaned; } ``` 这步看着简单,但很重要。Tika 提取出来的文本经常会有各种奇怪的控制字符、多余的空行,如果不清洗,后面的标题识别和段落切分都会受影响。 ## 小结 到这里,异步解析链路的入口和文本提取阶段就走完了。整个过程可以概括为: **消费 Kafka 消息 → 加载上下文 → 推进任务状态 → 下载原始文件 → Tika 提取文本 → 文本清洗** 拿到清洗后的纯文本之后,`parse()` 方法接下来要做的第一件事就是结构节点提取——这是整条解析流水线中最复杂的一环,下一篇单独展开。 --- **企业级项目导航**:⬅️ [[07-知识路由的功能|07-知识路由的功能]] | 01-Kafka 消费与文本内容解析 | ➡️ [[02-一定得用Kafka吗|02-一定得用Kafka吗]]